iT邦幫忙

2026 iThome 鐵人賽

DAY 14
0
Software Development

Kotlin 手刻 Ktor 從零開始系列 第 14

Kotlin 手刻 Ktor 從零開始 Day 14 升級 suspend,讓 Pipeline 支援協程

  • 分享至 

  • xImage
  •  

https://ithelp.ithome.com.tw/upload/images/20260822/20121948NIqCtgaIUJ.png

到第 13 篇為止,Relix 的 pipeline 是同步的,handler 和 middleware 都是 RelixCall.() -> RelixResponse,只要你開始做資料庫存取或呼叫外部 API,就會希望 handler 能用 suspend,把等待 I/O 的時間交給協程調度,而不是阻塞整條 thread

這篇要做的事,把 pipeline 的型別全面升級成 suspend,確保既有測試全部還是沒有問題的

為什麼要升級成 suspend ?

同步 handler 會讓等待 I/O 的期間也占著 executor thread,JDK HttpServer 的實際並行方式由設定的 Executor 決定,不論使用哪種 executor,框架內部支援 suspend 後,middleware 才能自然呼叫協程的 API

系列採取的策略是「先同步再非同步」,先把概念講清楚,再把型別升級,這樣學起來比較不會被協程細節淹沒

加 coroutines dependency

動程式碼之前,先把相依性補上,這也是系列第一次動到 module.yaml 的 dependencies

第 02 篇建出來的專案,module.yaml 從頭到尾只有一行

product: jvm/app

typealias 加 suspend 這件事本身不需要任何函式庫,suspend 是 Kotlin 的語言關鍵字,編譯器產生的 Continuation 也在 kotlin-stdlib 裡,但等一下 TestKit 和 JDK adapter 都要用 runBlocking 當 bridge,而 runBlocking 來自 kotlinx-coroutines-core,沒有它的話 import kotlinx.coroutines.runBlocking 那行會直接紅字

product: jvm/app

dependencies:
  - org.jetbrains.kotlinx:kotlinx-coroutines-core:1.11.0

這裡用的是當下的最新穩定版 1.11.0,你也可以到 MVN Repository 挑一個與自己 Kotlin Toolchain 相容的版本

加完先跑一次 ./kotlin build,確認相依性抓得到,再往下改型別

更新 typealias 型別

把原本在 Pipeline.ktRelixHandler.kt 同步的 typealias 全部加上 suspend

// 同步版(第 12-13 篇)
typealias RelixHandler = RelixCall.() -> RelixResponse
typealias Next = RelixCall.() -> RelixResponse
typealias RelixMiddleware = RelixCall.(next: Next) -> RelixResponse

// suspend 版(本篇升級)
typealias RelixHandler = suspend RelixCall.() -> RelixResponse
typealias Next = suspend RelixCall.() -> RelixResponse
typealias RelixMiddleware = suspend RelixCall.(next: Next) -> RelixResponse

每行只多了一個 suspend,型別的結構完全沒有變,差別只在這些 lambda 現在可以呼叫 suspend function 了

buildPipeline 不用動

Pipeline.kt 裡面的 buildPipeline,一個字都不用改

fun buildPipeline(handler: Next, middlewares: List<RelixMiddleware>): Next {
    return middlewares.foldRight(handler) { middleware, next ->
        { middleware(next) }
    }
}

因為 NextRelixMiddleware 的 typealias 已經改成 suspendfoldRight 裡面的 lambda 自動就是 suspend lambda。你也不需要在 buildPipeline 的函式簽名上加 suspend,因為它只是在組合函式,不是在執行它

這是 Kotlin 型別系統的好處,typealias 改了,依賴它的程式碼自動跟上

更新 RelixApplication.handle()

handle() 要變成 suspend function

class RelixApplication {

    suspend fun handle(call: RelixCall): RelixResponse {
        val terminal: Next = when (
            val result = router.match(call.request.method, call.request.path)
        ) {
            is MatchResult.Matched -> {
                call.pathParams = result.pathParams
                call.matchedRoute = result.route
                result.route.handler
            }
            is MatchResult.NotFound -> { { notFound() } }
            is MatchResult.MethodNotAllowed -> {
                { methodNotAllowed(result.allowedMethods) }
            }
        }

        val pipeline = buildPipeline(terminal, middlewares)
        val response = pipeline(call)

        return if (call.request.method == "HEAD") {
            response.copy(body = ByteArray(0))
        } else {
            response
        }
    }
}

唯一的差異是 fun handle 變成 suspend fun handle

如果你只改了 typealias,忘了改這裡,val response = pipeline(call) 這行會編譯失敗

Suspend function 'suspend fun invoke(p1: RelixCall): RelixResponse' can only be called from a coroutine or another suspend function.

原因是 pipeline 的型別是 Next,而 Next 已經加上 suspend 了,所以 pipeline(call) 變成一個 suspend 呼叫,它需要一個 suspend 的呼叫端,解法就是把 handle() 加上 suspend,函式裡面的邏輯一行都不用動

TestKit 升級,runBlocking bridge

handle() 變成 suspend 之後,TestKit 的 handleRequest() 不能直接呼叫它了,它會編譯錯誤,最簡單的 bridge 是用 runBlocking

import kotlinx.coroutines.runBlocking

class RelixTestKit(private val application: RelixApplication) {

    fun handleRequest(
        method: String = "GET",
        path: String = "/",
        headers: Map<String, List<String>> = emptyMap(),
        queryParameters: Map<String, List<String>> = emptyMap(),
        body: ByteArray = ByteArray(0),
    ): RelixResponse = runBlocking {
        val request = RelixRequest(
                   method = method,
                   path = path,
                   headers = headers,
                   queryParameters = queryParameters,
                   body = body,
        )
        val call = RelixCall(application, request)
        application.handle(call)
    }
}

handleRequest() 的簽名沒變,還是回傳 RelixResponse (不是 suspend),內部用 runBlocking 把 suspend handle 跑完,這代表走 TestKit 的既有測試,呼叫方式完全不用改

testApplication helper 也不需要改,因為它回傳的是 RelixTestKit,而 RelixTestKit 內部處理了 runBlocking

JDK HttpServer Adapter 的 bridge

跟 TestKit 一樣,JDK adapter 也需要 runBlocking

class JdkHttpServerAdapter(private val application: RelixApplication) {
    private lateinit var server: HttpServer
    var port: Int = 0
        private set

    fun start(port: Int = 8080) {
        try {
            server = HttpServer.create(InetSocketAddress(port), 0)
        } catch (e: java.net.BindException) {
            throw IllegalStateException(
                "Port $port is already in use. Try a different port or use port 0 for auto-assign.",
                e
            )
        }

        server.createContext("/") { exchange ->
            try {
                val request = exchange.toRelixRequest()
                val call = RelixCall(application, request)
                val response = runBlocking { application.handle(call) }
                exchange.writeResponse(response)
            } catch (e: Exception) {
                val body = "Internal Server Error".toByteArray()
                exchange.sendResponseHeaders(500, body.size.toLong())
                exchange.responseBody.use { it.write(body) }
            }
        }

        server.start()
        this.port = server.address.port
    }

    fun stop() {
        server.stop(0)
    }
}

關鍵是 runBlocking { application.handle(call) },JDK HttpServer 的 handler 是同步 callback,不在 coroutine context 裡,因此需要一個 bridge,這個版本會阻塞目前的 executor thread,suspend 本身不會把 JDK adapter 變成 non-blocking server

suspend 編譯後到底長什麼樣

只加一個關鍵字 pipeline 就升級了,看起來很簡單,但 compiler 背後做了不少事,Kotlin 會把 suspend function 編譯成 CPS (continuation passing style),函式簽名多一個隱藏的 Continuation 參數,函式內部也會被拆成狀態機,每個 suspend point 對應一個 state,它不是讓 thread 神奇消失,而是把「等 I/O」改寫成「先保存目前狀態,之後再接著跑」

// 你寫的
suspend fun handle(call: RelixCall): RelixResponse {
    val response = pipeline(call)        // suspend point
    return response
}

// compiler 在 bytecode 層面做的事(概念,不是真的程式碼)
fun handle(call: RelixCall, $cont: Continuation<RelixResponse>): Any {
    when ($cont.label) {
        0 -> {
            $cont.label = 1
            val result = pipeline(call, $cont)            // 可能立即回值,也可能回 COROUTINE_SUSPENDED
            if (result == COROUTINE_SUSPENDED) return result
            return result
        }
        1 -> return $cont.result                          // 被 resume 時從這裡接著跑
    }
}

runBlocking 的職責是在目前 thread 上跑完一段 coroutine,每當 continuation resume 時,它會把流程接回去,直到整條 chain 結束,對非 coroutine 的呼叫端來說,像 JDK HttpServer 的 handler thread、JUnit 的同步測試,這是一座必要的橋,它不是長期架構的最佳解,只是把同步入口接到 suspend pipeline 的務實做法

如果你想再深入,Roman Elizarov 的 KotlinConf talk「Deep Dive into Coroutines on JVM」把整個編譯流程用組譯碼一行一行示範,系列裡不深入這塊,但知道「suspend 不是 thread」這件事,足以讓你寫出正確的 middleware

先修好編譯失敗的既有測試

在跑測試之前,有兩個測試檔會先卡住,第 12 篇的 MiddlewareTest 和第 13 篇的 PipelineTest,它們的共通點是都沒有經過 TestKit,而是自己直接呼叫 middleware 或 pipeline

PipelineTest 是自己組 pipeline、自己呼叫

val pipeline = buildPipeline(handler, listOf(middlewareA, middlewareB))
val call = RelixCall(RelixApplication(), RelixRequest("GET", "/", emptyMap(), emptyMap(), ByteArray(0)))
pipeline(call)    // 編譯失敗

MiddlewareTest 則是直接把 middleware 當函式呼叫,自己傳一個 lambda 扮演 next

val response = requestIdMiddleware(call) {    // 編譯失敗
    seenInsideHandler = context.get<String>("traceId")
    ok("done")
}

兩個失敗的原因是同一個,pipeline 的型別是 NextrequestIdMiddleware 的型別是 RelixMiddleware,兩個 typealias 現在都帶著 suspend,而 @Test fun 是普通函式,錯誤訊息跟前面 handle() 那次一模一樣

修法也一樣,把呼叫包進 runBlocking

import kotlinx.coroutines.runBlocking

// PipelineTest
val response = runBlocking { pipeline(call) }

// MiddlewareTest
val response = runBlocking {
    requestIdMiddleware(call) {
        seenInsideHandler = context.get<String>("traceId")
        ok("done")
    }
}

MiddlewareTest 有四個測試,PipelineTest 有五個,每一個都有自己的呼叫點,要各包一次,但只有呼叫那一行要動,其它的寫法都不用改

middleware 和 handler 的定義也都不用動,不管是 val middlewareA: RelixMiddleware = { next -> ... } 還是第 12 篇那三個示範 middleware,這種沒有呼叫 suspend function 的普通 lambda 會被隱式轉成 suspend lambda,裡面的 next() 也因為身處 suspend lambda 而合法

直接碰 middleware、pipeline 或 handle() 的測試要補 runBlocking,走 RelixTestKit.handleRequest() 的測試一行都不用改,因為 runBlocking 已經在 TestKit 內部了

確認既有測試沒壞,再補新測試

編譯過了之後,先跑一次完整的測試套件,前面累積下來的測試就是這次升級最好的安全網,這次改的是型別不是行為,行為沒變,它們就應該要全部通過

確認舊的沒壞之後,再補三個專門給 suspend 版的測試,把洋蔥模型的 before/after 順序、短路、routing 各驗一次,確認它們在 suspend 之後還是原本的樣子

import kotlin.test.Test
import kotlin.test.assertEquals

class SuspendPipelineTest {

    @Test
    fun `suspend pipeline maintains before-after order`() {
        val trace = mutableListOf<String>()

        val app = RelixApplication()
        app.use { next ->
            trace += "A-before"
            val response = next()
            trace += "A-after"
            response
        }
        app.use { next ->
            trace += "B-before"
            val response = next()
            trace += "B-after"
            response
        }
        app.routing {
            get("/test") {
                trace += "handler"
                ok("done")
            }
        }

        val testKit = RelixTestKit(app)
        val response = testKit.handleRequest("GET", "/test")

        assertEquals(200, response.statusCode)
        assertEquals(
            listOf("A-before", "B-before", "handler", "B-after", "A-after"),
            trace,
        )
    }

    @Test
    fun `suspend middleware can short-circuit`() {
        val app = RelixApplication()
        app.use { _ ->
            RelixResponse(401, emptyMap(), "Unauthorized".toByteArray())
        }
        app.routing {
            get("/secret") { ok("secret data") }
        }

        val testKit = RelixTestKit(app)
        val response = testKit.handleRequest("GET", "/secret")

        assertEquals(401, response.statusCode)
    }

    @Test
    fun `existing routing tests still pass after suspend upgrade`() {
        val app = RelixApplication()
        app.routing {
            get("/hello") { ok("Hello!") }
            get("/users/{id}") { ok("User: ${pathParam("id")}") }
        }

        val testKit = RelixTestKit(app)

        assertEquals(200, testKit.handleRequest("GET", "/hello").statusCode)
        assertEquals("User: 42", testKit.handleRequest("GET", "/users/42").bodyAsText())
        assertEquals(404, testKit.handleRequest("GET", "/missing").statusCode)
    }
}

遷移清單

做這種橫跨多個檔案的型別升級,列一張清單會比較不容易漏,順序就是實際動手的順序

檔案 變更
module.yaml 加入 kotlinx-coroutines-core dependency
RelixHandler.kt typealias RelixHandler = suspend RelixCall.() -> RelixResponse
Pipeline.kt NextRelixMiddlewaresuspendbuildPipeline 不用動
RelixApplication.kt handle()suspend
RelixTestKit.kt handleRequest() 內部用 runBlocking
JdkHttpServerAdapter.kt handle() 呼叫處加 runBlocking
MiddlewareTest.kt middleware 呼叫各包一次 runBlocking
PipelineTest.kt pipeline(call) 呼叫各包一次 runBlocking
SuspendPipelineTest.kt 三個新的 suspend 版本測試

module.yaml 要先加,runBlocking 沒有相依性就用不了

中間三個 src/ 檔案是型別升級的本體,接下來四個 bridge 現場其實是同一件事,凡是從非 suspend 的地方呼叫 suspend 的東西,都要有一座橋

RelixApplicationTest.kt 不在清單裡,因為它走的是 TestKit,runBlocking 已經在 handleRequest() 內部了

常見陷阱與設計取捨

runBlocking 不是不好嗎 ? 為什麼到處用 ?

runBlocking 會阻塞目前的 thread,因此不能放在 Netty 之類的 event loop 上,這裡只把它當成 JDK 同步 callback 與 suspend pipeline 之間的橋接,實際吞吐量仍受 executor 設定、阻塞式工作與 thread 數量影響,不能直接推論成與同步版完全相同

如果之後換成 Netty 或 CIO 引擎,就要把 runBlocking 換成 coroutine-native 的處理方式,但那是引擎層的事,framework 內部的 suspend pipeline 不用改

handler 裡不寫 suspend 的話會怎樣 ?

完全沒問題,普通的 lambda (沒有呼叫任何 suspend function) 可以隱式轉成 suspend lambda,所以之前寫的所有 handler { ok("Hello!") } 不用改,compiler 會自動處理

為什麼先同步再 suspend,不直接從 suspend 開始 ?

兩個原因,第一,如果同時學 pipeline 概念和 coroutine,注意力會分散,先用同步版把洋蔥模型、fold 串鏈、短路行為都搞懂,再加上 suspend,一次只改一個變數,第二,同步版的測試是 suspend 版的安全網,如果 suspend 升級出了問題,你有一整套同步版的測試可以對照


小結

升級 suspend 改的東西不多,三個 typealias 加 suspendhandle()suspend、TestKit 和 Adapter 用 runBlocking bridge。buildPipeline 的組合邏輯完全不用動。走 TestKit 的既有測試不用改寫法,只有直接呼叫 middleware 或 pipeline 的 MiddlewareTestPipelineTest 要補 runBlocking。先同步再 suspend 的漸進策略讓這次升級變成一個低風險的重構,而不是一次大改寫


下一篇

下一篇開始加第一個內建 middleware,Logging,用它驗證 pipeline 的 before/after 行為,也讓框架在開發時更好用


參考資料


同步刊登於 Blog

圖片來源:AI 產生


上一篇
Kotlin 手刻 Ktor 從零開始 Day 13 實作同步版 Pipeline,把 Middleware 串成鏈
下一篇
Kotlin 手刻 Ktor 從零開始 Day 15 內建 Middleware,Logging
系列文
Kotlin 手刻 Ktor 從零開始18
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言